🌎 rns.kin.earth git 🌍
Commit 99c428a9f5406560a8ae7247630b2947eba8cc8d
Parents : cb9071b
Author : Mark Qvist <bc7291552be7a58f361522990465165c>
Signature : T66BB85Valid, signed by author
Date : 2026-08-27T02:22:55+02:00
Added dataplane egress control
Changes
3 files changed, 110 insertions(+), 5 deletions(-)
Diff
diff --git a/RNS/Interfaces/BackboneInterface.py b/RNS/Interfaces/BackboneInterface.py
index c41483f8..ce476540 100644
--- a/RNS/Interfaces/BackboneInterface.py
+++ b/RNS/Interfaces/BackboneInterface.py
@@ -60,6 +60,13 @@ class BackboneInterface(Interface):
DP_IC_RCVBUF = 32768
DP_IC_IF_HEADROOM = 32
DP_IC_PENALTY = 1.5
+ DP_EC_INTERVAL = 1.0
+ DP_EC_MID_WM = 128*1024
+ DP_EC_HIGH_WM = 4*1024*1024
+ DP_EC_STALL_TICKS = 3
+ DP_EC_MAX_ETA = 10.0
+ DP_EC_RELEASE_ETA = 5.0 # Gate release hysterisis
+ DP_EC_DEAD_TIME = 12.0
epoll = None
listener_filenos = {}
@@ -67,7 +74,9 @@ class BackboneInterface(Interface):
_job_active = False
_ic_job_active = False
+ _ec_job_active = False
_job_lock = threading.Lock()
+ _ec_job_lock = threading.Lock()
_ic_job_lock = threading.Lock()
_dp_ic_lock = threading.Lock()
_dp_ic_snapshot = 0
@@ -253,6 +262,7 @@ class BackboneInterface(Interface):
def start():
if not BackboneInterface._job_active: threading.Thread(target=BackboneInterface.__job, daemon=True).start()
if not BackboneInterface._ic_job_active: threading.Thread(target=BackboneInterface.__dp_ic_job, daemon=True).start()
+ if not BackboneInterface._ec_job_active: threading.Thread(target=BackboneInterface.__dp_ec_job, daemon=True).start()
@staticmethod
def ensure_epoll():
@@ -398,6 +408,86 @@ class BackboneInterface(Interface):
RNS.log(f"Available ingress budget {RNS.prettyspeed(avail*8)}, hard-allocated {round(alloc*100.0, 2)}% ({RNS.prettyspeed(avail*alloc*8)}), handled in {RNS.prettyshorttime(taken, compact=True, tight=True)}", RNS.LOG_DEBUG)
RNS.log(f"Holding for {RNS.prettyshorttime(hold, compact=True, tight=True)}, throttling handled in {RNS.prettyshorttime(taken, compact=True, tight=True)} at depth {q_depth}", RNS.LOG_DEBUG)
+ @staticmethod
+ def _dp_ec_evaluate(interface, now):
+ tb = interface.transmit_buffer
+ drained = tb._tx_sent - interface._dp_ec_prev_sent
+ interface._dp_ec_prev_sent = tb._tx_sent
+
+ sendable = tb.sendable
+ buffered = len(tb)
+
+ if buffered == 0 or sendable == 0:
+ interface._dp_ec_zero_ticks = 0
+ interface._dp_ec_last_drain = now
+ interface.tx_stalled = False
+ return False
+
+ if now - interface._dp_ec_last_drain >= BackboneInterface.DP_EC_DEAD_TIME:
+ RNS.log(f"No egress control drain progress for {RNS.prettyshorttime(BackboneInterface.DP_EC_DEAD_TIME, compact=True)} on {interface}, tearing down", RNS.LOG_WARNING)
+ try:
+ if hasattr(interface, "socket") and interface.socket:
+ fileno = interface.socket.fileno()
+ BackboneInterface.deregister_fileno(fileno)
+ if fileno in BackboneInterface.spawned_interface_filenos: BackboneInterface.spawned_interface_filenos.pop(fileno)
+ try: interface.socket.close()
+ except Exception as e: RNS.log(f"Egress control could not close socket for {interface}: {e}", RNS.LOG_ERROR)
+ except Exception as e: RNS.log(f"Egress control cleanup error for {interface}: {e}", RNS.LOG_ERROR)
+
+ interface.receive(b"")
+ return True
+
+ if drained > 0:
+ interface._dp_ec_last_drain = now
+ interface._dp_ec_zero_ticks = 0
+ drain_rate = drained / BackboneInterface.DP_EC_INTERVAL
+ clear_eta = buffered / drain_rate if drain_rate > 0 else float("inf")
+
+ if buffered > BackboneInterface.DP_EC_MID_WM and clear_eta > BackboneInterface.DP_EC_MAX_ETA:
+ if not interface.tx_stalled: RNS.log(f"Egress control drain ETA of {RNS.prettyshorttime(clear_eta, compact=True)} exceeds maximum on {interface}, gating outbound", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
+ interface.tx_stalled = True
+
+ elif clear_eta < BackboneInterface.DP_EC_RELEASE_ETA or buffered <= BackboneInterface.DP_EC_MID_WM:
+ if interface.tx_stalled: RNS.log(f"Egress control drain recovered on {interface}, resuming outbound", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
+ interface.tx_stalled = False
+
+ else:
+ # No drain progress on this tick
+ if buffered > BackboneInterface.DP_EC_MID_WM:
+ interface._dp_ec_zero_ticks += 1
+ if interface._dp_ec_zero_ticks >= BackboneInterface.DP_EC_STALL_TICKS:
+ if not interface.tx_stalled: RNS.log(f"No egress control drain progress on {interface}, gating outbound", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
+ interface.tx_stalled = True
+ else:
+ interface._dp_ec_zero_ticks = 0
+ interface.tx_stalled = False
+
+ return False
+
+ @staticmethod
+ def __dp_ec_job():
+ with BackboneInterface._ec_job_lock:
+ if BackboneInterface._ec_job_active: return
+ else:
+ BackboneInterface._ec_job_active = True
+ RNS.log(f"Started dataplane egress control", RNS.LOG_DEBUG)
+ try:
+ while True:
+ time.sleep(BackboneInterface.DP_EC_INTERVAL)
+ now = time.time()
+ try: interfaces = list(BackboneInterface.spawned_interface_filenos.values())
+ except RuntimeError:
+ RNS.log(f"Deferring egress control evaluation due to error: {e}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
+ continue
+
+ for interface in interfaces:
+ if interface.detached: continue
+ BackboneInterface._dp_ec_evaluate(interface, now)
+
+ except Exception as e:
+ RNS.log(f"BackboneInterface egress control error: {e}", RNS.LOG_ERROR)
+ RNS.trace_exception(e)
+
@staticmethod
def __dp_ic_job():
with BackboneInterface._ic_job_lock:
@@ -416,7 +506,6 @@ class BackboneInterface(Interface):
with BackboneInterface._dp_ic_lock:
st = time.time()
interfaces = list(BackboneInterface.spawned_interface_filenos.values())
- # active = [iface for iface in interfaces if iface.dp_ingress_bytes > 0]
if q_depth > BackboneInterface.DP_IC_MID_WM:
producers = [iface for iface in interfaces if iface.dp_ingress_packets > 0]
@@ -941,8 +1030,12 @@ class BackboneClientInterface(Interface):
def process_outgoing(self, data):
if self.online and not self.detached:
try:
- self.transmit_buffer.append(bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]))
- BackboneInterface.tx_ready(self)
+ frame = bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG])
+ if self.tx_stalled or not self.transmit_buffer.append(frame, self.tx_hwm):
+ self.tx_drops += 1
+ self.tx_dropped_bytes += len(frame)
+ RNS.log(f"Egress control dropping outbound frame of {RNS.prettysize(len(frame))} on {self}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
+ else: BackboneInterface.tx_ready(self)
except Exception as e:
RNS.log("Exception occurred while transmitting via "+str(self)+", tearing down interface", RNS.LOG_ERROR)
diff --git a/RNS/Interfaces/Interface.py b/RNS/Interfaces/Interface.py
index 7823cbe9..ca301f86 100755
--- a/RNS/Interfaces/Interface.py
+++ b/RNS/Interfaces/Interface.py
@@ -163,6 +163,14 @@ class Interface:
self.dp_ingress_hold = None
self.dp_ingress_gated = False
+ self.tx_hwm = 4*1024*1024
+ self.tx_stalled = False
+ self.tx_drops = 0
+ self.tx_dropped_bytes = 0
+ self._dp_ec_prev_sent = 0
+ self._dp_ec_zero_ticks = 0
+ self._dp_ec_last_drain = time.time()
+
self.reports_phy_stats = False
self.r_stat_rssi = None
self.r_stat_snr = None
diff --git a/RNS/Interfaces/LocalInterface.py b/RNS/Interfaces/LocalInterface.py
index 330bd345..cfc5eb2b 100644
--- a/RNS/Interfaces/LocalInterface.py
+++ b/RNS/Interfaces/LocalInterface.py
@@ -218,8 +218,12 @@ class LocalClientInterface(Interface):
if self.online:
try:
if self.epoll_backend:
- self.transmit_buffer.append(bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG]))
- BackboneInterface.tx_ready(self)
+ frame = bytes([HDLC.FLAG])+HDLC.escape(data)+bytes([HDLC.FLAG])
+ if self.tx_stalled or not self.transmit_buffer.append(frame, self.tx_hwm):
+ self.tx_drops += 1
+ self.tx_dropped_bytes += len(frame)
+ RNS.log(f"Egress control dropping outbound frame of {RNS.prettysize(len(frame))} on {self}", RNS.LOG_DEBUG) if RNS.sl(RNS.LOG_DEBUG) else None
+ else: BackboneInterface.tx_ready(self)
else:
self.writing = True
Served by rngit 1.5.2 - Generated in 0.05s